package net.xdclass.dwd;

import net.xdclass.util.KafkaUtil;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;

/**
 * @Author: wangzhan
 * @Description: TODO
 * @DateTime: 2025/2/15 15:41
 **/
public class TestFlink {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.setParallelism(1);

        DataStream<String> ds = env.socketTextStream("127.0.0.1", 8888);

//        FlinkKafkaConsumer<String> kafkaConsumer = KafkaUtil.getKafkaConsumer(SOURCE_TOPIC, GROUP_ID);

//        DataStreamSource<String> ds = env.addSource(kafkaConsumer);

        ds.print();
        env.execute();
    }
}
